package org.example.flatmap;

import org.apache.commons.lang3.StringUtils;
import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.util.Collector;

/**
 * 处理
 *
 * @author wlj
 */
public class MyFlatMap implements FlatMapFunction<String, String> {

    private static final long serialVersionUID = 7723150291726141635L;

    @Override
    public void flatMap(String s, Collector<String> collector) {
        // s 是kafka消息，collector是搜集之后返回的数据
        if (StringUtils.isNotBlank(s)) {

            // 针对消息进行一系列的处理加工
            // TODO

            collector.collect(s);
        }
    }
}
